--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
Commit 241b29cc92b15907bc0c17cfb7a41a17669e50ac
Parents : ca02fa5
Author : Mark Qvist <bc7291552be7a58f361522990465165c>
Signature : T66BB85Valid, signed by author
Date : 2026-07-16T02:50:10+02:00
Yes, it will run on potato, then
Changes
2 files changed, 107 insertions(+), 20 deletions(-)
Diff
diff --git a/LXMF/LXMRouter.py b/LXMF/LXMRouter.py
index 83e741d..7f17541 100644
--- a/LXMF/LXMRouter.py
+++ b/LXMF/LXMRouter.py
@@ -53,6 +53,9 @@ class LXMRouter:
PROPAGATION_COST_FLEX = 3
PROPAGATION_COST = 16
PROPAGATION_LIMIT = 256
+ SEQUENTIAL_VALIDATION = True
+ STATIC_SEQUENTIAL = False
+ MAX_INBOUND_SYNCS = 3
SYNC_LIMIT = PROPAGATION_LIMIT*40
DELIVERY_LIMIT = 1000
@@ -76,6 +79,11 @@ class LXMRouter:
PR_ALL_MESSAGES = 0x00
+ OFFER_UNKNOWN = 0x00
+ OFFER_ACCEPTED = 0x01
+ OFFER_TRANSFERRING = 0x02
+ OFFER_VALIDATING = 0x03
+
DUPLICATE_SIGNAL = "lxmf_duplicate"
STATS_GET_PATH = "/pn/get/stats"
@@ -86,12 +94,13 @@ class LXMRouter:
### Developer-facing API ##############################
#######################################################
- def __init__(self, identity=None, storagepath=None, autopeer=AUTOPEER, autopeer_maxdepth=None,
+ def __init__(self, identity=None, storagepath=None, name=None, autopeer=AUTOPEER, autopeer_maxdepth=None,
propagation_limit=PROPAGATION_LIMIT, delivery_limit=DELIVERY_LIMIT, sync_limit=SYNC_LIMIT,
enforce_ratchets=False, enforce_stamps=False, static_peers = [], max_peers=None,
from_static_only=False, sync_strategy=LXMPeer.STRATEGY_PERSISTENT,
propagation_cost=PROPAGATION_COST, propagation_cost_flexibility=PROPAGATION_COST_FLEX,
- peering_cost=PEERING_COST, max_peering_cost=MAX_PEERING_COST, name=None):
+ peering_cost=PEERING_COST, max_peering_cost=MAX_PEERING_COST, max_inbound_syncs=MAX_INBOUND_SYNCS,
+ sequential_validation=SEQUENTIAL_VALIDATION, static_sequential=STATIC_SEQUENTIAL):
random.seed(os.urandom(10))
@@ -131,6 +140,9 @@ class LXMRouter:
self.information_storage_limit = None
self.propagation_per_transfer_limit = propagation_limit
self.propagation_per_sync_limit = sync_limit
+ self.propagation_sequential_validation = sequential_validation
+ self.propagation_static_peer_sequential = static_sequential
+ self.propagation_max_inbound_syncs = max_inbound_syncs
self.delivery_per_transfer_limit = delivery_limit
self.propagation_stamp_cost = propagation_cost
self.propagation_stamp_cost_flexibility = propagation_cost_flexibility
@@ -153,6 +165,7 @@ class LXMRouter:
self.propagation_transfer_max_messages = None
self.prioritise_rotating_unreachable_peers = False
self.active_propagation_links = []
+ self.accepted_offer_links = {}
self.validated_peer_links = {}
self.locally_delivered_transient_ids = {}
self.locally_processed_transient_ids = {}
@@ -165,6 +178,7 @@ class LXMRouter:
self.cost_file_lock = threading.Lock()
self.ticket_file_lock = threading.Lock()
self.stamp_gen_lock = threading.Lock()
+ self.accepted_offer_links_lock = threading.Lock()
self.exit_handler_running = False
if identity == None:
@@ -174,6 +188,8 @@ class LXMRouter:
self.propagation_destination = RNS.Destination(self.identity, RNS.Destination.IN, RNS.Destination.SINGLE, APP_NAME, "propagation")
self.propagation_destination.set_default_app_data(self.get_propagation_node_app_data)
self.control_destination = None
+ self.validating_pn_stamps_from = {}
+ self.sequential_validation_lock = threading.Lock()
self.client_propagation_messages_received = 0
self.client_propagation_messages_served = 0
self.unpeered_propagation_incoming = 0
@@ -934,6 +950,17 @@ class LXMRouter:
for link in inactive_links:
self.active_propagation_links.remove(link)
link.teardown()
+
+ active_link_ids = []
+ inactive_offers = []
+ for link in self.active_propagation_links: active_link_ids.append(link.link_id)
+ with self.accepted_offer_links_lock:
+ for link_id in self.accepted_offer_links:
+ if not link_id in active_link_ids: inactive_offers.append(link_id)
+
+ for link_id in inactive_offers:
+ RNS.log(f"Cleaning inbound sync link accounting for link {RNS.prettyhexrep(link_id)} since link is no longer active", RNS.LOG_DEBUG) # TODO: Remove at some point
+ self.accepted_offer_links.pop(link_id)
except Exception as e:
RNS.log("An error occurred while cleaning inbound propagation links. The contained exception was: "+str(e), RNS.LOG_ERROR)
@@ -1343,10 +1370,8 @@ class LXMRouter:
self.propagation_destination.deregister_request_handler(LXMRouter.UNPEER_REQUEST_PATH)
for link in self.active_propagation_links:
try:
- if link.status == RNS.Link.ACTIVE:
- link.teardown()
- except Exception as e:
- RNS.log("Error while tearing down propagation link: {e}", RNS.LOG_ERROR)
+ if link.status == RNS.Link.ACTIVE: link.teardown()
+ except Exception as e: RNS.log("Error while tearing down propagation link: {e}", RNS.LOG_ERROR)
RNS.log("Persisting LXMF state data to storage...", RNS.LOG_NOTICE)
self.flush_queues()
@@ -1886,8 +1911,7 @@ class LXMRouter:
if resource.status == RNS.Resource.COMPLETE:
ratchet_id = None
# Set ratchet ID to link ID if available
- if resource.link and hasattr(resource.link, "link_id"):
- ratchet_id = resource.link.link_id
+ if resource.link and hasattr(resource.link, "link_id"): ratchet_id = resource.link.link_id
phy_stats = {"rssi": resource.link.rssi, "snr": resource.link.snr, "q": resource.link.q}
self.lxmf_delivery(resource.data.read(), resource.link.type, phy_stats=phy_stats, ratchet_id=ratchet_id, method=LXMessage.DIRECT)
@@ -2092,6 +2116,16 @@ class LXMRouter:
link.set_resource_concluded_callback(self.propagation_resource_concluded)
self.active_propagation_links.append(link)
+ @property
+ def propagation_resources_transferring(self):
+ count = 0
+ with self.accepted_offer_links_lock:
+ for link_id in self.accepted_offer_links:
+ if self.accepted_offer_links[link_id] > self.OFFER_ACCEPTED:
+ count += 1
+
+ return count
+
def propagation_resource_advertised(self, resource):
if self.from_static_only:
remote_identity = resource.link.get_remote_identity()
@@ -2111,7 +2145,13 @@ class LXMRouter:
if limit != None and size > limit:
RNS.log(f"Rejecting {RNS.prettysize(size)} incoming propagation resource, since it exceeds the limit of {RNS.prettysize(limit)}", RNS.LOG_DEBUG)
return False
+
else:
+ with self.accepted_offer_links_lock:
+ if resource.link.link_id in self.accepted_offer_links:
+ ri_str = RNS.prettyhexrep(resource.link.get_remote_identity().hash) if resource.link.get_remote_identity() else 'unknown peer'
+ RNS.log(f"Sync offer for {ri_str} started transferring", RNS.LOG_DEBUG) # TODO: Remove at some point
+ self.accepted_offer_links[resource.link.link_id] = self.OFFER_TRANSFERRING
return True
def propagation_packet(self, data, packet):
@@ -2147,13 +2187,24 @@ class LXMRouter:
RNS.log("The contained exception was: "+str(e), RNS.LOG_ERROR)
def offer_request(self, path, data, request_id, link_id, remote_identity, requested_at):
- if remote_identity == None:
- return LXMPeer.ERROR_NO_IDENTITY
+ if remote_identity == None: return LXMPeer.ERROR_NO_IDENTITY
else:
remote_destination = RNS.Destination(remote_identity, RNS.Destination.OUT, RNS.Destination.SINGLE, APP_NAME, "propagation")
remote_hash = remote_destination.hash
remote_str = RNS.prettyhexrep(remote_hash)
+ bypass_sequential = not self.propagation_static_peer_sequential and remote_hash in self.static_peers
+ if not bypass_sequential and self.propagation_sequential_validation and len(self.validating_pn_stamps_from) > 0:
+ RNS.log(f"Propagation offer from node {remote_str} postponed, already validating {len(self.validating_pn_stamps_from)} PN stamp batches", RNS.LOG_NOTICE)
+ # for rh in self.validating_pn_stamps_from:
+ # RNS.log(f"Validating from {RNS.prettyhexrep(rh)} for {RNS.prettytime(time.time()-self.validating_pn_stamps_from[rh])}")
+ return LXMPeer.ERROR_THROTTLED
+
+ resources_transferring = self.propagation_resources_transferring
+ if not bypass_sequential and self.propagation_max_inbound_syncs and resources_transferring >= self.propagation_max_inbound_syncs:
+ RNS.log(f"Propagation offer from node {remote_str} postponed, already receiving {resources_transferring} sync resource{'s' if resources_transferring != 1 else ''}", RNS.LOG_NOTICE)
+ return LXMPeer.ERROR_THROTTLED
+
if remote_hash in self.throttled_peers:
throttle_remaining = self.throttled_peers[remote_hash]-time.time()
if throttle_remaining > 0:
@@ -2189,9 +2240,16 @@ class LXMRouter:
for transient_id in transient_ids:
if not transient_id in self.propagation_entries: wanted_ids.append(transient_id)
- if len(wanted_ids) == 0: return False
- elif len(wanted_ids) == len(transient_ids): return True
- else: return wanted_ids
+ if len(wanted_ids) == 0:
+ RNS.log(f"No wanted messages in offer from {RNS.prettyhexrep(remote_hash)}", RNS.LOG_DEBUG)
+ return False
+ elif len(wanted_ids) == len(transient_ids):
+ RNS.log(f"Accepted all {len(wanted_ids)} offered message{'s' if len(wanted_ids) != 1 else ''} from {RNS.prettyhexrep(remote_hash)}", RNS.LOG_DEBUG)
+ return True
+ else:
+ RNS.log(f"Accepted {len(wanted_ids)} offered message{'s' if len(wanted_ids) != 1 else ''} from {RNS.prettyhexrep(remote_hash)}", RNS.LOG_DEBUG)
+ with self.accepted_offer_links_lock: self.accepted_offer_links[link_id] = self.OFFER_ACCEPTED
+ return wanted_ids
except Exception as e:
RNS.log("Error occurred while generating response for sync request, the contained exception was: "+str(e), RNS.LOG_DEBUG)
@@ -2199,6 +2257,7 @@ class LXMRouter:
return None
def propagation_resource_concluded(self, resource):
+ with self.accepted_offer_links_lock: self.accepted_offer_links[resource.link.link_id] = self.OFFER_VALIDATING
if resource.status == RNS.Resource.COMPLETE:
try:
data = msgpack.unpackb(resource.data.read())
@@ -2252,12 +2311,33 @@ class LXMRouter:
ms = "" if len(messages) == 1 else "s"
RNS.log(f"Received {len(messages)} message{ms} from {remote_str}, validating stamps...", RNS.LOG_VERBOSE)
- min_accepted_cost = max(0, self.propagation_stamp_cost-self.propagation_stamp_cost_flexibility)
- validated_messages = LXStamper.validate_pn_stamps(messages, min_accepted_cost)
- invalid_stamps = len(messages)-len(validated_messages)
- ms = "" if invalid_stamps == 1 else "s"
- if len(validated_messages) == len(messages): RNS.log(f"All message stamps validated from {remote_str}", RNS.LOG_VERBOSE)
- else: RNS.log(f"Transfer from {remote_str} contained {invalid_stamps} invalid stamp{ms}", RNS.LOG_WARNING)
+ with self.sequential_validation_lock:
+ RNS.log(f"Adding validation job accounting entry for {RNS.prettyhexrep(remote_hash)}", RNS.LOG_DEBUG) # TODO: Remove at some point
+ self.validating_pn_stamps_from[remote_hash] = time.time()
+ try:
+ min_accepted_cost = max(0, self.propagation_stamp_cost-self.propagation_stamp_cost_flexibility)
+ validated_messages = LXStamper.validate_pn_stamps(messages, min_accepted_cost)
+ invalid_stamps = len(messages)-len(validated_messages)
+ ms = "" if invalid_stamps == 1 else "s"
+ if len(validated_messages) == len(messages): RNS.log(f"All message stamps validated from {remote_str}", RNS.LOG_VERBOSE)
+ else: RNS.log(f"Transfer from {remote_str} contained {invalid_stamps} invalid stamp{ms}", RNS.LOG_WARNING)
+
+ except Exception as e:
+ RNS.log(f"Error while validating received propagation message stamps: {e}", RNS.LOG_ERROR)
+ RNS.trace_exception(e)
+ return
+
+ finally:
+ with self.sequential_validation_lock:
+ RNS.log(f"Cleaning validation job accounting entry for {RNS.prettyhexrep(remote_hash)}", RNS.LOG_DEBUG) # TODO: Remove at some point
+ try: self.validating_pn_stamps_from.pop(remote_hash)
+ except Exception as e: RNS.log(f"Failed to remove PN stamp validation job from sequential tracking: {e}", RNS.LOG_ERROR)
+
+ with self.accepted_offer_links_lock:
+ if resource.link.link_id in self.accepted_offer_links:
+ ri_str = RNS.prettyhexrep(resource.link.get_remote_identity().hash) if resource.link.get_remote_identity() else 'unknown peer'
+ RNS.log(f"Cleaning inbound sync link accounting for {ri_str}", RNS.LOG_DEBUG) # TODO: Remove at some point
+ self.accepted_offer_links.pop(resource.link.link_id)
for validated_entry in validated_messages:
transient_id = validated_entry[0]
@@ -2296,6 +2376,12 @@ class LXMRouter:
RNS.log("Error while unpacking received propagation resource", RNS.LOG_DEBUG)
RNS.trace_exception(e)
+ with self.accepted_offer_links_lock:
+ if resource.link.link_id in self.accepted_offer_links:
+ ri_str = RNS.prettyhexrep(resource.link.get_remote_identity().hash) if resource.link.get_remote_identity() else 'unknown peer'
+ RNS.log(f"Cleaning inbound sync link accounting for {ri_str} on resource failure", RNS.LOG_DEBUG) # TODO: Remove at some point
+ self.accepted_offer_links.pop(resource.link.link_id)
+
def enqueue_peer_distribution(self, transient_id, from_peer):
self.peer_distribution_queue.append([transient_id, from_peer])
diff --git a/LXMF/Utilities/lxmd.py b/LXMF/Utilities/lxmd.py
index 1b13acb..fa2e658 100644
--- a/LXMF/Utilities/lxmd.py
+++ b/LXMF/Utilities/lxmd.py
@@ -1105,7 +1105,8 @@ autopeer_maxdepth = 6
# propagation sync transfers will be accepted.
# Once this number is reached, nodes offering
# messages will receive a throttle response to
-# indicate that they can retry later.
+# indicate that they can retry later. On a slow
+# system, it's a good idea to change this to 1.
# max_inbound_syncs = 3
──────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────